Conversation
There was a problem hiding this comment.
Code Review
This pull request updates the Grafana dashboard's circuit breaker state query to map states to numeric values and configures Kafka consumers to use ErrorHandlingDeserializer to prevent deserialization errors from blocking the consumer. The reviewer provided two valuable suggestions: first, using StringJsonMessageConverter in KafkaConfig.java to safely handle messages without type headers from external systems or CLIs; second, including the forced_open state in the Grafana PromQL expression to prevent false-positive normal statuses when a circuit breaker is manually forced open.
| config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); | ||
| config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); | ||
| config.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class.getName()); | ||
| config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName()); | ||
|
|
||
| config.put(JsonDeserializer.TRUSTED_PACKAGES, "com.catchtable.notification.event,java.lang.String,java.lang.Object"); | ||
| // 헤더 미포함 메시지(외부 시스템/CLI) 대비 fallback. 헤더가 있으면 우선 적용된다. | ||
| config.put(JsonDeserializer.USE_TYPE_INFO_HEADERS, true); |
There was a problem hiding this comment.
헤더 미포함 메시지 처리를 위한 StringJsonMessageConverter 도입 제안
현재 설정은 JsonDeserializer를 Consumer 레벨에서 직접 사용하고 있습니다. 주석에 언급된 것처럼 외부 시스템이나 CLI 등에서 헤더(__TypeId__) 없이 전송된 메시지를 처리하려 할 때, JsonDeserializer는 역직렬화할 대상 클래스 타입을 알 수 없어 IllegalArgumentException을 발생시키며 실패하게 됩니다. (단일 default type을 지정하려 해도, 현재 컨슈머 그룹 내에서 여러 종류의 이벤트 클래스(ReservationConfirmedEvent, VacancyEvent 등)를 처리하고 있어 설정이 불가능합니다.)
이 문제를 해결하는 가장 표준적이고 견고한 방법은 **StringJsonMessageConverter**를 사용하는 것입니다.
- Consumer Factory에서는 단순히
StringDeserializer를 사용하여 메시지를 문자열로 읽어옵니다. ConcurrentKafkaListenerContainerFactory에StringJsonMessageConverter를 등록하면, Spring Kafka가@KafkaListener메서드의 파라미터 타입(예:ReservationConfirmedEvent)을 기반으로 JSON 문자열을 해당 객체로 자동 역직렬화해 줍니다.- 이 방식을 사용하면
__TypeId__헤더 유무와 상관없이 모든 외부 메시지를 안전하게 수신할 수 있습니다.
아래는 kafkaListenerContainerFactory 변경 예시입니다.
@Bean
public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory(KafkaTemplate<String, Object> kafkaTemplate) {
ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
// StringJsonMessageConverter 등록 (메서드 파라미터 기반 자동 역직렬화)
factory.setRecordMessageConverter(new StringJsonMessageConverter());
DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(
kafkaTemplate,
(record, ex) -> new TopicPartition(record.topic() + "-dlt", record.partition())
);
factory.setCommonErrorHandler(new DefaultErrorHandler(recoverer, new FixedBackOff(2000L, 3L)));
return factory;
}| config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); | |
| config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); | |
| config.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class.getName()); | |
| config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName()); | |
| config.put(JsonDeserializer.TRUSTED_PACKAGES, "com.catchtable.notification.event,java.lang.String,java.lang.Object"); | |
| // 헤더 미포함 메시지(외부 시스템/CLI) 대비 fallback. 헤더가 있으면 우선 적용된다. | |
| config.put(JsonDeserializer.USE_TYPE_INFO_HEADERS, true); | |
| config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); | |
| config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class); | |
| config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, StringDeserializer.class.getName()); |
| "description": "closed=정상, half_open=복구시도, open=차단됨. AI API 등 외부 의존성 보호.", | ||
| "type": "stat", | ||
| "targets": [ { "expr": "resilience4j_circuitbreaker_state == 1", "legendFormat": "{{name}} ({{state}})", "refId": "A", "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" } } ] | ||
| "targets": [ { "expr": "sum by (name) (resilience4j_circuitbreaker_state{state=\"closed\"} * 0) + sum by (name) (resilience4j_circuitbreaker_state{state=\"half_open\"} * 1) + sum by (name) (resilience4j_circuitbreaker_state{state=\"open\"} * 2)", "legendFormat": "{{name}}", "refId": "A", "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" } } ] |
There was a problem hiding this comment.
forced_open 및 disabled 상태 누락으로 인한 오탐지 위험
현재 PromQL 식은 closed (0), half_open (1), open (2) 상태만 고려하여 합산하고 있습니다.
하지만 Resilience4j에는 forced_open 및 disabled 상태도 존재합니다. 만약 운영 중에 긴급 대응이나 테스트를 위해 서킷 브레이커를 forced_open 상태로 강제 전환할 경우, closed/half_open/open 메트릭은 모두 0이 되고 forced_open 메트릭만 1이 됩니다.
이 경우 현재 식의 결과는 0이 되어 대시보드에는 정상 상태인 CLOSED (초록색)로 표시되는 오탐지(Misleading)가 발생할 수 있습니다.
따라서 forced_open 상태도 open과 동일하게 2로 매핑할 수 있도록 정규식(=~"open|forced_open")을 사용하는 것을 권장합니다.
| "targets": [ { "expr": "sum by (name) (resilience4j_circuitbreaker_state{state=\"closed\"} * 0) + sum by (name) (resilience4j_circuitbreaker_state{state=\"half_open\"} * 1) + sum by (name) (resilience4j_circuitbreaker_state{state=\"open\"} * 2)", "legendFormat": "{{name}}", "refId": "A", "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" } } ] | |
| "targets": [ { "expr": "sum by (name) (resilience4j_circuitbreaker_state{state=\"closed\"} * 0) + sum by (name) (resilience4j_circuitbreaker_state{state=\"half_open\"} * 1) + sum by (name) (resilience4j_circuitbreaker_state{state=~\"open|forced_open\"} * 2)", "legendFormat": "{{name}}", "refId": "A", "datasource": { "type": "prometheus", "uid": "DS_PROMETHEUS" } } ] |
📢 기능 설명
필요시 실행결과 스크린샷 첨부
연결된 issue
연결된 issue를 자동을 닫기 위해 아래 {이슈넘버}를 입력해주세요.
close #{이슈넘버}
✅ 체크리스트